🎖️GitЯра🎖️
Node / meshtastic / Meshtastic-Android / files / core / network / src / commonMain / kotlin / org / meshtastic / core / network / radio / TcpRadioTransport.kt
Displaying Raw • Download
core/network/src/commonMain/kotlin/org/meshtastic/core/network/radio/TcpRadioTransport.kt e9d09a3380eca7860f1b8dfbd1939d8415fc0e33 (e9d09a33) Text, 8.11 KB
T8b949e/*
* Copyright (c) 2026 Meshtastic LLC
*
* This program is free software: you can redistribute it and/or modify
* it under the terms of the GNU General Public License as published by
* the Free Software Foundation, either version 3 of the License, or
* (at your option) any later version.
*
* This program is distributed in the hope that it will be useful,
* but WITHOUT ANY WARRANTY; without even the implied warranty of
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
* GNU General Public License for more details.
*
* You should have received a copy of the GNU General Public License
* along with this program. If not, see <https://www.gnu.org/licenses/>.
*/
Tff7b72package T7ee787org.meshtastic.core.network.radio
Tff7b72import T7ee787co.touchlab.kermit.Logger
Tff7b72import T7ee787kotlinx.atomicfu.atomic
Tff7b72import T7ee787kotlinx.atomicfu.locks.SynchronizedObject
Tff7b72import T7ee787kotlinx.atomicfu.locks.synchronized
Tff7b72import T7ee787kotlinx.coroutines.CoroutineScope
Tff7b72import T7ee787kotlinx.coroutines.Job
Tff7b72import T7ee787kotlinx.coroutines.joinAll
Tff7b72import T7ee787kotlinx.coroutines.withTimeoutOrNull
Tff7b72import T7ee787org.meshtastic.core.common.util.handledLaunch
Tff7b72import T7ee787org.meshtastic.core.di.CoroutineDispatchers
Tff7b72import T7ee787org.meshtastic.core.network.transport.StreamFrameCodec
Tff7b72import T7ee787org.meshtastic.core.network.transport.TcpTransport
Tff7b72import T7ee787org.meshtastic.core.repository.RadioTransport
Tff7b72import T7ee787org.meshtastic.core.repository.RadioTransportCallback
Tff7b72import T7ee787kotlin.time.Duration.Companion.seconds
T8b949e/** Minimal TCP connection surface used by [TcpRadioTransport] to coordinate admission and teardown. */
Tff7b72internal Tff7b72interface T56d364TcpRadioConnection Tb4b4b4{
Tff7b72val Te6edf3isConnectedTb4b4b4: Tffa657Boolean
Tff7b72fun Td2a8ffstartTb4b4b4(Te6edf3addressTb4b4b4: Tffa657StringTb4b4b4)
Tff7b72fun Td2a8ffstopTb4b4b4(Tb4b4b4)
Tff7b72suspend Tff7b72fun Td2a8ffsendPacketTb4b4b4(Te6edf3payloadTb4b4b4: Te6edf3ByteArrayTb4b4b4)
Tff7b72suspend Tff7b72fun Td2a8ffsendHeartbeatTb4b4b4(Tb4b4b4)
Tb4b4b4}
Tff7b72private Tff7b72class T56d364TcpRadioConnectionImplTb4b4b4(Tff7b72private Tff7b72val Te6edf3delegateTb4b4b4: Te6edf3TcpTransportTb4b4b4) Tb4b4b4: Te6edf3TcpRadioConnection Tb4b4b4{
Tff7b72override Tff7b72val Te6edf3isConnectedTb4b4b4: Tffa657Boolean
Tff7b72getTb4b4b4(Tb4b4b4) Tff7b72= Te6edf3delegateTb4b4b4.Te6edf3isConnected
Tff7b72override Tff7b72fun Td2a8ffstartTb4b4b4(Te6edf3addressTb4b4b4: Tffa657StringTb4b4b4) Tff7b72= Te6edf3delegateTb4b4b4.Te6edf3startTb4b4b4(Te6edf3addressTb4b4b4)
Tff7b72override Tff7b72fun Td2a8ffstopTb4b4b4(Tb4b4b4) Tff7b72= Te6edf3delegateTb4b4b4.Te6edf3stopTb4b4b4(Tb4b4b4)
Tff7b72override Tff7b72suspend Tff7b72fun Td2a8ffsendPacketTb4b4b4(Te6edf3payloadTb4b4b4: Te6edf3ByteArrayTb4b4b4) Tff7b72= Te6edf3delegateTb4b4b4.Te6edf3sendPacketTb4b4b4(Te6edf3payloadTb4b4b4)
Tff7b72override Tff7b72suspend Tff7b72fun Td2a8ffsendHeartbeatTb4b4b4(Tb4b4b4) Tff7b72= Te6edf3delegateTb4b4b4.Te6edf3sendHeartbeatTb4b4b4(Tb4b4b4)
Tb4b4b4}
T8b949e/**
* TCP radio transport — thin adapter over the shared [TcpTransport] from `core:network`.
*
* Implements [RadioTransport] directly via composition over [TcpTransport], delegating send/receive to the transport
* and calling [RadioTransportCallback] for lifecycle events. This avoids the previous inheritance from
* [StreamTransport] which created a dead [StreamFrameCodec] and required overriding `sendBytes` as a no-op.
*/
Tff7b72open Tff7b72class T56d364TcpRadioTransport
Te6edf3internal Tff7b72constructorTb4b4b4(
Tff7b72private Tff7b72val Te6edf3callbackTb4b4b4: Te6edf3RadioTransportCallbackTb4b4b4,
Tff7b72private Tff7b72val Te6edf3scopeTb4b4b4: Te6edf3CoroutineScopeTb4b4b4,
Tff7b72private Tff7b72val Te6edf3addressTb4b4b4: Tffa657StringTb4b4b4,
Te6edf3connectionFactoryTb4b4b4: Tb4b4b4(Te6edf3TcpTransportTb4b4b4.Te6edf3ListenerTb4b4b4) Tff7b72-Tff7b72> Te6edf3TcpRadioConnectionTb4b4b4,
Tb4b4b4) Tb4b4b4: Te6edf3RadioTransport Tb4b4b4{
Tff7b72constructorTb4b4b4(
Te6edf3callbackTb4b4b4: Te6edf3RadioTransportCallbackTb4b4b4,
Te6edf3scopeTb4b4b4: Te6edf3CoroutineScopeTb4b4b4,
Te6edf3dispatchersTb4b4b4: Te6edf3CoroutineDispatchersTb4b4b4,
Te6edf3addressTb4b4b4: Tffa657StringTb4b4b4,
Tb4b4b4) Tb4b4b4: Tff7b72thisTb4b4b4(
Te6edf3callback Tff7b72= Te6edf3callbackTb4b4b4,
Te6edf3scope Tff7b72= Te6edf3scopeTb4b4b4,
Te6edf3address Tff7b72= Te6edf3addressTb4b4b4,
Te6edf3connectionFactory Tff7b72= Tb4b4b4{ Te6edf3listener Tff7b72-Tff7b72>
Te6edf3TcpRadioConnectionImplTb4b4b4(
Te6edf3TcpTransportTb4b4b4(
Te6edf3dispatchers Tff7b72= Te6edf3dispatchersTb4b4b4,
Te6edf3scope Tff7b72= Te6edf3scopeTb4b4b4,
Te6edf3listener Tff7b72= Te6edf3listenerTb4b4b4,
Te6edf3logTag Tff7b72= Ta5d6ff"Ta5d6ffTcpRadioTransport[Tffd700$Te6edf3addressTa5d6ff]Ta5d6ff"Tb4b4b4,
Tb4b4b4)Tb4b4b4,
Tb4b4b4)
Tb4b4b4}Tb4b4b4,
Tb4b4b4)
Tff7b72companion Tff7b72object Tb4b4b4{
Tff7b72const Tff7b72val Te6edf3SERVICE_PORT Tff7b72= Te6edf3StreamFrameCodecTb4b4b4.Te6edf3DEFAULT_TCP_PORT
Tff7b72internal Tff7b72val Te6edf3OPERATION_TIMEOUT Tff7b72= T79c0ff1T79c0ff0.Te6edf3seconds
Tb4b4b4}
Tff7b72private Tff7b72val Te6edf3lifecycle Tff7b72= Te6edf3TransportLifecycleGateTb4b4b4(Ta5d6ff"Ta5d6ffTCPTa5d6ff"Tb4b4b4, Te6edf3operationDrainTimeout Tff7b72= Te6edf3OPERATION_TIMEOUT Tff7b72* T79c0ff2Tb4b4b4)
Tff7b72private Tff7b72val Te6edf3transportStopped Tff7b72= Te6edf3atomicTb4b4b4(Tff7b72falseTb4b4b4)
Tff7b72private Tff7b72val Te6edf3operationJobsLock Tff7b72= Te6edf3SynchronizedObjectTb4b4b4(Tb4b4b4)
Tff7b72private Tff7b72val Te6edf3operationJobs Tff7b72= Te6edf3mutableSetOfTff7b72<Te6edf3JobTff7b72>Tb4b4b4(Tb4b4b4)
Tff7b72private Tff7b72val Te6edf3transport Tff7b72=
Te6edf3connectionFactoryTb4b4b4(
Tff7b72object Tb4b4b4: T56d364TcpTransportTb4b4b4.Te6edf3Listener Tb4b4b4{
Tff7b72override Tff7b72fun Td2a8ffonConnectedTb4b4b4(Tb4b4b4) Tb4b4b4{
Te6edf3lifecycleTb4b4b4.Te6edf3runIfOpen Tb4b4b4{ Te6edf3callbackTb4b4b4.Te6edf3onConnectTb4b4b4(Tb4b4b4) Tb4b4b4}
Tb4b4b4}
Tff7b72override Tff7b72fun Td2a8ffonDisconnectedTb4b4b4(Tb4b4b4) Tb4b4b4{
T8b949e// close() first closes admission, suppressing the transient callback caused by transport.stop().
Te6edf3lifecycleTb4b4b4.Te6edf3runIfOpen Tb4b4b4{ Te6edf3callbackTb4b4b4.Te6edf3onDisconnectTb4b4b4(Te6edf3isPermanent Tff7b72= Tff7b72falseTb4b4b4) Tb4b4b4}
Tb4b4b4}
Tff7b72override Tff7b72fun Td2a8ffonPacketReceivedTb4b4b4(Te6edf3bytesTb4b4b4: Te6edf3ByteArrayTb4b4b4) Tb4b4b4{
Te6edf3lifecycleTb4b4b4.Te6edf3runIfOpen Tb4b4b4{ Te6edf3callbackTb4b4b4.Te6edf3handleFromRadioTb4b4b4(Te6edf3bytesTb4b4b4) Tb4b4b4}
Tb4b4b4}
Tb4b4b4}Tb4b4b4,
Tb4b4b4)
Tff7b72override Tff7b72fun Td2a8ffstartTb4b4b4(Tb4b4b4) Tb4b4b4{
Te6edf3lifecycleTb4b4b4.Te6edf3runIfOpen Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3transportStoppedTb4b4b4.Te6edf3valueTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{ Ta5d6ff"Ta5d6ff[Tffd700$Te6edf3addressTa5d6ff] Ignoring start on a stopped TCP transport; a fresh transport is requiredTa5d6ff" Tb4b4b4}
Tb4b4b4} Tff7b72else Tb4b4b4{
Te6edf3transportTb4b4b4.Te6edf3startTb4b4b4(Te6edf3addressTb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tff7b72override Tff7b72suspend Tff7b72fun Td2a8ffcloseTb4b4b4(Tb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ff[Tffd700$Te6edf3addressTa5d6ff] Closing TCP transportTa5d6ff" Tb4b4b4}
Tff7b72val Te6edf3completed Tff7b72=
Te6edf3lifecycleTb4b4b4.Te6edf3closeTb4b4b4(
Te6edf3teardown Tff7b72= Tb4b4b4{
Te6edf3stopTransportTb4b4b4(Tb4b4b4)
Te6edf3cancelOutstandingOperationsTb4b4b4(Tb4b4b4)
Tb4b4b4}Tb4b4b4,
Tb4b4b4)
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3completedTb4b4b4) Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{ Ta5d6ff"Ta5d6ff[Tffd700$Te6edf3addressTa5d6ff] TCP teardown did not complete within its lifecycle boundsTa5d6ff" Tb4b4b4}
T8b949e// Do NOT emit onDisconnect(isPermanent = true) here. The explicit-disconnect signal is the service layer's
T8b949e// responsibility (SharedRadioInterfaceService.stopTransportLocked); emitting it here causes a double-disconnect
T8b949e// and prevents the auto-reconnect loop from owning its transient lifecycle.
Tb4b4b4}
Tff7b72override Tff7b72fun Td2a8ffkeepAliveTb4b4b4(Tb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ff[Tffd700$Te6edf3addressTa5d6ff] TCP keepAliveTa5d6ff" Tb4b4b4}
Te6edf3launchConnectionOperationTb4b4b4(Ta5d6ff"Ta5d6ffheartbeatTa5d6ff"Tb4b4b4) Tb4b4b4{ Te6edf3transportTb4b4b4.Te6edf3sendHeartbeatTb4b4b4(Tb4b4b4) Tb4b4b4}
Tb4b4b4}
Tff7b72override Tff7b72fun Td2a8ffhandleSendToRadioTb4b4b4(Te6edf3pTb4b4b4: Te6edf3ByteArrayTb4b4b4)Tb4b4b4: Tffa657Boolean Tff7b72=
Te6edf3launchConnectionOperationTb4b4b4(Ta5d6ff"Ta5d6ffsendTa5d6ff"Tb4b4b4) Tb4b4b4{ Te6edf3transportTb4b4b4.Te6edf3sendPacketTb4b4b4(Te6edf3pTb4b4b4) Tb4b4b4}
T8b949e/**
* Schedules one operation while the current transport is connected.
*
* The Boolean reports lifecycle admission and successful scheduling only; the operation runs asynchronously and may
* still fail or time out after this method returns `true`.
*/
Tff7b72private Tff7b72fun Td2a8fflaunchConnectionOperationTb4b4b4(Te6edf3operationTb4b4b4: Tffa657StringTb4b4b4, Te6edf3blockTb4b4b4: Te6edf3suspend Tb4b4b4(Tb4b4b4) Tff7b72-Tff7b72> Tffa657UnitTb4b4b4)Tb4b4b4: Tffa657Boolean Tb4b4b4{
Tff7b72val Te6edf3lease Tff7b72= Te6edf3lifecycleTb4b4b4.Te6edf3tryAcquireTb4b4b4(Tb4b4b4) Tff7b72?: Tff7b72return Tff7b72false
Tff7b72return Tff7b72if Tb4b4b4(Tff7b72!Te6edf3transportTb4b4b4.Te6edf3isConnectedTb4b4b4) Tb4b4b4{
Te6edf3leaseTb4b4b4.Te6edf3releaseTb4b4b4(Tb4b4b4)
Tff7b72false
Tb4b4b4} Tff7b72else Tb4b4b4{
Tff7b72var Te6edf3completionRegistered Tff7b72= Tff7b72false
Tff7b72try Tb4b4b4{
Tff7b72val Te6edf3job Tff7b72=
Te6edf3scopeTb4b4b4.Te6edf3handledLaunch Tb4b4b4{
Tff7b72val Te6edf3completed Tff7b72= Te6edf3withTimeoutOrNullTb4b4b4(Te6edf3OPERATION_TIMEOUTTb4b4b4) Tb4b4b4{ Te6edf3blockTb4b4b4(Tb4b4b4) Tb4b4b4} Tff7b72!Tff7b72= Tff7b72null
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3completedTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{ Ta5d6ff"Ta5d6ff[Tffd700$Te6edf3addressTa5d6ff] TCP Tffd700$Te6edf3operationTa5d6ff timed out after Tffd700$Te6edf3OPERATION_TIMEOUTTa5d6ff" Tb4b4b4}
T8b949e// Cancellation may leave a framed write partially emitted. Stopping this one-shot
T8b949e// transport forces the service reconnect path to create a fresh transport before another
T8b949e// send.
Te6edf3stopTransportTb4b4b4(Tb4b4b4)
Tb4b4b4}
Tb4b4b4}
Te6edf3synchronizedTb4b4b4(Te6edf3operationJobsLockTb4b4b4) Tb4b4b4{ Te6edf3operationJobs Tff7b72+Tff7b72= Te6edf3job Tb4b4b4}
Te6edf3jobTb4b4b4.Te6edf3invokeOnCompletion Tb4b4b4{
Te6edf3synchronizedTb4b4b4(Te6edf3operationJobsLockTb4b4b4) Tb4b4b4{ Te6edf3operationJobs Tff7b72-Tff7b72= Te6edf3job Tb4b4b4}
Te6edf3leaseTb4b4b4.Te6edf3releaseTb4b4b4(Tb4b4b4)
Tb4b4b4}
Te6edf3completionRegistered Tff7b72= Tff7b72true
Tff7b72!Te6edf3jobTb4b4b4.Te6edf3isCancelled
Tb4b4b4} Tff7b72finally Tb4b4b4{
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3completionRegisteredTb4b4b4) Te6edf3leaseTb4b4b4.Te6edf3releaseTb4b4b4(Tb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffcancelOutstandingOperationsTb4b4b4(Tb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3jobs Tff7b72= Te6edf3synchronizedTb4b4b4(Te6edf3operationJobsLockTb4b4b4) Tb4b4b4{ Te6edf3operationJobsTb4b4b4.Te6edf3toListTb4b4b4(Tb4b4b4) Tb4b4b4}
Te6edf3jobsTb4b4b4.Te6edf3forEach Tb4b4b4{ Tffa657itTb4b4b4.Te6edf3cancelTb4b4b4(Tb4b4b4) Tb4b4b4}
Tff7b72val Te6edf3joined Tff7b72= Te6edf3withTimeoutOrNullTb4b4b4(Te6edf3OPERATION_TIMEOUTTb4b4b4) Tb4b4b4{ Te6edf3jobsTb4b4b4.Te6edf3joinAllTb4b4b4(Tb4b4b4) Tb4b4b4} Tff7b72!Tff7b72= Tff7b72null
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3joinedTb4b4b4) Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{ Ta5d6ff"Ta5d6ff[Tffd700$Te6edf3addressTa5d6ff] TCP operation jobs did not stop after transport teardownTa5d6ff" Tb4b4b4}
Tb4b4b4}
Tff7b72private Tff7b72fun Td2a8ffstopTransportTb4b4b4(Tb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3transportStoppedTb4b4b4.Te6edf3compareAndSetTb4b4b4(Te6edf3expect Tff7b72= Tff7b72falseTb4b4b4, Te6edf3update Tff7b72= Tff7b72trueTb4b4b4)Tb4b4b4) Te6edf3transportTb4b4b4.Te6edf3stopTb4b4b4(Tb4b4b4)
Tb4b4b4}
Tb4b4b4}
Served by rngit 1.5.2 - Generated in 0.05s